🎖️GitЯра🎖️
Node / meshtastic / Meshtastic-Android / files / core / data / src / commonMain / kotlin / org / meshtastic / core / data / manager / PacketHandlerImpl.kt
Displaying Raw • Download
core/data/src/commonMain/kotlin/org/meshtastic/core/data/manager/PacketHandlerImpl.kt 65a1f5ce4d0e4b8eafcd0bbfccafd41dc6ee888a (65a1f5ce) Text, 35.09 KB
T8b949e/*
* Copyright (c) 2026 Meshtastic LLC
*
* This program is free software: you can redistribute it and/or modify
* it under the terms of the GNU General Public License as published by
* the Free Software Foundation, either version 3 of the License, or
* (at your option) any later version.
*
* This program is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
* GNU General Public License for more details.
*
* You should have received a copy of the GNU General Public License
* along with this program. If not, see <https://www.gnu.org/licenses/>.
*/
Tff7b72package T7ee787org.meshtastic.core.data.manager
Tff7b72import T7ee787co.touchlab.kermit.Logger
Tff7b72import T7ee787kotlinx.atomicfu.atomic
Tff7b72import T7ee787kotlinx.coroutines.CancellationException
Tff7b72import T7ee787kotlinx.coroutines.CompletableDeferred
Tff7b72import T7ee787kotlinx.coroutines.CoroutineStart
Tff7b72import T7ee787kotlinx.coroutines.Deferred
Tff7b72import T7ee787kotlinx.coroutines.Job
Tff7b72import T7ee787kotlinx.coroutines.NonCancellable
Tff7b72import T7ee787kotlinx.coroutines.TimeoutCancellationException
Tff7b72import T7ee787kotlinx.coroutines.delay
Tff7b72import T7ee787kotlinx.coroutines.isActive
Tff7b72import T7ee787kotlinx.coroutines.selects.select
Tff7b72import T7ee787kotlinx.coroutines.sync.Mutex
Tff7b72import T7ee787kotlinx.coroutines.sync.withLock
Tff7b72import T7ee787kotlinx.coroutines.withContext
Tff7b72import T7ee787kotlinx.coroutines.withTimeout
Tff7b72import T7ee787kotlinx.coroutines.withTimeoutOrNull
Tff7b72import T7ee787org.koin.core.annotation.Single
Tff7b72import T7ee787org.meshtastic.core.common.di.ServiceScope
Tff7b72import T7ee787org.meshtastic.core.common.util.handledLaunch
Tff7b72import T7ee787org.meshtastic.core.common.util.nowMillis
Tff7b72import T7ee787org.meshtastic.core.database.dao.shouldApplyOutgoingQueueStatus
Tff7b72import T7ee787org.meshtastic.core.model.ConnectionState
Tff7b72import T7ee787org.meshtastic.core.model.MeshLog
Tff7b72import T7ee787org.meshtastic.core.model.MessageStatus
Tff7b72import T7ee787org.meshtastic.core.model.RadioNotConnectedException
Tff7b72import T7ee787org.meshtastic.core.model.util.toOneLineString
Tff7b72import T7ee787org.meshtastic.core.model.util.toPIIString
Tff7b72import T7ee787org.meshtastic.core.repository.AwaitedSendResult
Tff7b72import T7ee787org.meshtastic.core.repository.AwaitedSendStatus
Tff7b72import T7ee787org.meshtastic.core.repository.ConnectionStateProvider
Tff7b72import T7ee787org.meshtastic.core.repository.MeshLogRepository
Tff7b72import T7ee787org.meshtastic.core.repository.PacketHandler
Tff7b72import T7ee787org.meshtastic.core.repository.PacketRepository
Tff7b72import T7ee787org.meshtastic.core.repository.PersistedPacketId
Tff7b72import T7ee787org.meshtastic.core.repository.PersistedReactionId
Tff7b72import T7ee787org.meshtastic.core.repository.RadioInterfaceService
Tff7b72import T7ee787org.meshtastic.proto.FromRadio
Tff7b72import T7ee787org.meshtastic.proto.MeshPacket
Tff7b72import T7ee787org.meshtastic.proto.QueueStatus
Tff7b72import T7ee787org.meshtastic.proto.Routing
Tff7b72import T7ee787org.meshtastic.proto.ToRadio
Tff7b72import T7ee787kotlin.time.Duration
Tff7b72import T7ee787kotlin.time.Duration.Companion.milliseconds
Tff7b72import T7ee787kotlin.time.Duration.Companion.minutes
Tff7b72import T7ee787kotlin.time.Duration.Companion.seconds
Tff7b72import T7ee787kotlin.uuid.Uuid
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooManyFunctionsTa5d6ff"Tb4b4b4)
Tf0883e@Single
Tff7b72class T56d364PacketHandlerImplTb4b4b4(
Tff7b72private Tff7b72val Te6edf3packetRepositoryTb4b4b4: Te6edf3LazyTff7b72<Te6edf3PacketRepositoryTff7b72>Tb4b4b4,
Tff7b72private Tff7b72val Te6edf3radioInterfaceServiceTb4b4b4: Te6edf3RadioInterfaceServiceTb4b4b4,
Tff7b72private Tff7b72val Te6edf3meshLogRepositoryTb4b4b4: Te6edf3LazyTff7b72<Te6edf3MeshLogRepositoryTff7b72>Tb4b4b4,
Tff7b72private Tff7b72val Te6edf3connectionStateProviderTb4b4b4: Te6edf3ConnectionStateProviderTb4b4b4,
Tff7b72private Tff7b72val Te6edf3scopeTb4b4b4: Te6edf3ServiceScopeTb4b4b4,
Tb4b4b4) Tb4b4b4: Te6edf3PacketHandler Tb4b4b4{
Tff7b72companion Tff7b72object Tb4b4b4{
Tff7b72internal Tff7b72val Te6edf3RESPONSE_TIMEOUT Tff7b72= T79c0ff5.Te6edf3seconds
T8b949e/** Routing acknowledgements traverse the mesh and need a larger budget than the local queue response. */
Tff7b72internal Tff7b72val Te6edf3ROUTING_RESPONSE_TIMEOUT Tff7b72= T79c0ff3T79c0ff0.Te6edf3seconds
Tff7b72private Tff7b72val Te6edf3PERSISTED_STATUS_WAIT Tff7b72= T79c0ff1.Te6edf3seconds
Tff7b72private Tff7b72val Te6edf3PERSISTED_STATUS_SHUTDOWN_WAIT Tff7b72= T79c0ff1T79c0ff5T79c0ff0.Te6edf3milliseconds
Tff7b72internal Tff7b72val Te6edf3PERSISTED_STATUS_RETRY_DELAY Tff7b72= T79c0ff1T79c0ff0T79c0ff0.Te6edf3milliseconds
T8b949e/**
* Grace period after which a sent packet still [MessageStatus.ENROUTE] is stamped [Routing.Error.TIMEOUT]
* (retryable) instead of showing as sending forever. Generous — well past the radio's retransmit window — and
* matches iOS's sendAckTimeout so both apps time out alike.
*/
Tff7b72internal Tff7b72val Te6edf3SEND_ACK_TIMEOUT Tff7b72= T79c0ff5.Te6edf3minutes
T8b949e/**
* Minimum re-arm delay on reconnect: the firmware's phone-queue backlog may still deliver the missing ACK/NAK
* just after the config handshake, so give it a moment before stamping a timeout.
*/
Tff7b72internal Tff7b72val Te6edf3REARM_GRACE Tff7b72= T79c0ff3T79c0ff0.Te6edf3seconds
T8b949e/**
* Firmware-internal `ErrorCode` (MeshTypes.h `ERRNO_SHOULD_RELEASE`) leaked into `QueueStatus.res`: "no error,
* but the packet should still be released". Firmware 2.8+ returns it for self-addressed packets, which are
* delivered through the synchronous local loopback instead of the TX queue — a success. Note it numerically
* collides with `Routing.Error.PKI_UNKNOWN_PUBKEY` (35); `QueueStatus.res` carries ErrorCode semantics, not
* Routing.Error.
*/
Tff7b72private Tff7b72const Tff7b72val Te6edf3ERRNO_SHOULD_RELEASE Tff7b72= T79c0ff3T79c0ff5
Tb4b4b4}
Tff7b72private Tff7b72var Te6edf3queueJobTb4b4b4: Te6edf3Job? Tff7b72= Tff7b72null
Tff7b72private Tff7b72var Te6edf3queueGeneration Tff7b72= T79c0ff0L
T8b949e// Lock-order invariant: acquire queueMutex before responseMutex whenever both are needed. Never acquire
T8b949e// responseMutex first on a path that can later need queueMutex, or admission and teardown can deadlock.
Tff7b72private Tff7b72val Te6edf3queueMutex Tff7b72= Te6edf3MutexTb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3queuedPackets Tff7b72= Te6edf3mutableListOfTff7b72<Te6edf3QueuedPacketTff7b72>Tb4b4b4(Tb4b4b4)
T8b949e// Marks the current queue generation as drained. The next admission clears it so reconnect recovery can start a
T8b949e// fresh worker; stale workers retain their generation and cannot replace that worker.
Tff7b72private Tff7b72var Te6edf3queueStopped Tff7b72= Tff7b72false
Tff7b72private Tff7b72val Te6edf3responseMutex Tff7b72= Te6edf3MutexTb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72data Tff7b72class T56d364QueuedPacketTb4b4b4(
Tff7b72val Te6edf3packetTb4b4b4: Te6edf3MeshPacketTb4b4b4,
Tff7b72val Te6edf3pendingTb4b4b4: Te6edf3PendingResponseTb4b4b4,
Tff7b72val Te6edf3expectedConnectionVersionTb4b4b4: Tffa657Long?Tb4b4b4,
Tb4b4b4)
Tff7b72private Tff7b72enum Tff7b72class T56d364StatusPersistence Tb4b4b4{
Te6edf3DATA_PACKETTb4b4b4,
Te6edf3REACTIONTb4b4b4,
Tb4b4b4}
Tff7b72private Tff7b72data Tff7b72class T56d364PacketStatusTargetTb4b4b4(Tff7b72val Te6edf3packetTb4b4b4: Te6edf3MeshPacketTb4b4b4, Tff7b72val Te6edf3persistenceTb4b4b4: Te6edf3StatusPersistence?Tb4b4b4)
Tff7b72private Tff7b72sealed Tff7b72interface T56d364PersistedStatusTarget Tb4b4b4{
Tff7b72data Tff7b72class T56d364DataPacketTb4b4b4(Tff7b72val Te6edf3idTb4b4b4: Te6edf3PersistedPacketIdTb4b4b4) Tb4b4b4: Te6edf3PersistedStatusTarget
Tff7b72data Tff7b72class T56d364ReactionTb4b4b4(Tff7b72val Te6edf3idTb4b4b4: Te6edf3PersistedReactionIdTb4b4b4) Tb4b4b4: Te6edf3PersistedStatusTarget
Tb4b4b4}
Tff7b72private Tff7b72class T56d364PendingResponseTb4b4b4(Tff7b72val Te6edf3packetTb4b4b4: Te6edf3MeshPacketTb4b4b4, Tff7b72val Te6edf3persistenceTb4b4b4: Te6edf3StatusPersistence?Tb4b4b4, Te6edf3awaitRoutingTb4b4b4: Tffa657BooleanTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3queueDeferred Tff7b72= Te6edf3CompletableDeferredTff7b72<Te6edf3AwaitedSendStatusTff7b72>Tb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3routingDeferred Tff7b72= Te6edf3CompletableDeferredTff7b72<Te6edf3AwaitedSendStatusTff7b72>Tb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3takeIf Tb4b4b4{ Te6edf3awaitRouting Tb4b4b4}
Tff7b72val Te6edf3dispatchAccepted Tff7b72= Te6edf3CompletableDeferredTff7b72<Tffa657BooleanTff7b72>Tb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3dispatchResolved Tff7b72= Te6edf3atomicTb4b4b4(Tff7b72falseTb4b4b4)
Tff7b72private Tff7b72val Te6edf3dispatchDepartureEpoch Tff7b72= Te6edf3atomicTff7b72<Tffa657Long?Tff7b72>Tb4b4b4(Tff7b72nullTb4b4b4)
Tff7b72private Tff7b72val Te6edf3queueStatus Tff7b72= Te6edf3atomicTff7b72<Te6edf3AwaitedSendStatus?Tff7b72>Tb4b4b4(Tff7b72nullTb4b4b4)
Tff7b72private Tff7b72val Te6edf3routingStatus Tff7b72= Te6edf3atomicTff7b72<Te6edf3AwaitedSendStatus?Tff7b72>Tb4b4b4(Tff7b72nullTb4b4b4)
Tff7b72val Te6edf3wasDispatchedTb4b4b4: Tffa657Boolean
Tff7b72getTb4b4b4(Tb4b4b4) Tff7b72= Te6edf3dispatchDepartureEpochTb4b4b4.Te6edf3value Tff7b72!Tff7b72= Tff7b72null
Tff7b72val Te6edf3departureEpochAtDispatchTb4b4b4: Tffa657Long?
Tff7b72getTb4b4b4(Tb4b4b4) Tff7b72= Te6edf3dispatchDepartureEpochTb4b4b4.Te6edf3value
Tff7b72val Te6edf3terminalStatusTb4b4b4: Te6edf3AwaitedSendStatus?
Tff7b72getTb4b4b4(Tb4b4b4) Tff7b72= Te6edf3routingStatusTb4b4b4.Te6edf3value Tff7b72?: Te6edf3queueStatusTb4b4b4.Te6edf3value
Tff7b72val Te6edf3isTerminalTb4b4b4: Tffa657Boolean
Tff7b72getTb4b4b4(Tb4b4b4) Tff7b72= Te6edf3queueDeferredTb4b4b4.Te6edf3isCompleted Tff7b72&Tff7b72& Tb4b4b4(Te6edf3routingDeferredTff7b72?.Te6edf3isCompleted Tff7b72!Tff7b72= Tff7b72falseTb4b4b4)
Tff7b72fun Td2a8ffrecordDispatchTb4b4b4(Te6edf3acceptedTb4b4b4: Tffa657BooleanTb4b4b4, Te6edf3departureEpochTb4b4b4: Tffa657LongTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3dispatchResolvedTb4b4b4.Te6edf3compareAndSetTb4b4b4(Te6edf3expect Tff7b72= Tff7b72falseTb4b4b4, Te6edf3update Tff7b72= Tff7b72trueTb4b4b4)Tb4b4b4) Tff7b72return
Tff7b72if Tb4b4b4(Te6edf3acceptedTb4b4b4) Te6edf3dispatchDepartureEpochTb4b4b4.Te6edf3value Tff7b72= Te6edf3departureEpoch
Te6edf3dispatchAcceptedTb4b4b4.Te6edf3completeTb4b4b4(Te6edf3acceptedTb4b4b4)
Tb4b4b4}
Tff7b72fun Td2a8ffcompleteQueueTb4b4b4(Te6edf3statusTb4b4b4: Te6edf3AwaitedSendStatusTb4b4b4, Te6edf3routingTerminalTb4b4b4: Tffa657Boolean Tff7b72= Te6edf3status Tff7b72!Tff7b72= Te6edf3AwaitedSendStatusTb4b4b4.Te6edf3ACCEPTEDTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3queueDeferredTb4b4b4.Te6edf3completeTb4b4b4(Te6edf3statusTb4b4b4)Tb4b4b4) Te6edf3queueStatusTb4b4b4.Te6edf3compareAndSetTb4b4b4(Tff7b72nullTb4b4b4, Te6edf3statusTb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3routingTerminalTb4b4b4) Te6edf3completeRoutingTb4b4b4(Te6edf3statusTb4b4b4)
Tb4b4b4}
Tff7b72fun Td2a8ffcompleteRoutingTb4b4b4(Te6edf3statusTb4b4b4: Te6edf3AwaitedSendStatusTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3deferred Tff7b72= Te6edf3routingDeferred Tff7b72?: Tff7b72return
Tff7b72if Tb4b4b4(Te6edf3deferredTb4b4b4.Te6edf3completeTb4b4b4(Te6edf3statusTb4b4b4)Tb4b4b4) Te6edf3routingStatusTb4b4b4.Te6edf3compareAndSetTb4b4b4(Tff7b72nullTb4b4b4, Te6edf3statusTb4b4b4)
Tb4b4b4}
Tff7b72fun Td2a8ffcompleteAllTb4b4b4(Te6edf3statusTb4b4b4: Te6edf3AwaitedSendStatusTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3dispatchResolvedTb4b4b4.Te6edf3compareAndSetTb4b4b4(Te6edf3expect Tff7b72= Tff7b72falseTb4b4b4, Te6edf3update Tff7b72= Tff7b72trueTb4b4b4)Tb4b4b4) Te6edf3dispatchAcceptedTb4b4b4.Te6edf3completeTb4b4b4(Tff7b72falseTb4b4b4)
Te6edf3completeQueueTb4b4b4(Te6edf3statusTb4b4b4, Te6edf3routingTerminal Tff7b72= Tff7b72trueTb4b4b4)
Tb4b4b4}
Tff7b72fun Td2a8ffsendFailureStatusTb4b4b4(Tb4b4b4)Tb4b4b4: Te6edf3AwaitedSendStatus Tff7b72=
Tff7b72if Tb4b4b4(Te6edf3wasDispatchedTb4b4b4) Te6edf3AwaitedSendStatusTb4b4b4.Te6edf3TRANSPORT_STOPPED Tff7b72else Te6edf3AwaitedSendStatusTb4b4b4.Te6edf3SEND_FAILED
Tb4b4b4}
Tff7b72private Tff7b72sealed Tff7b72interface T56d364QueueAdmission Tb4b4b4{
Tff7b72data Tff7b72class T56d364AdmittedTb4b4b4(Tff7b72val Te6edf3pendingTb4b4b4: Te6edf3PendingResponseTb4b4b4) Tb4b4b4: Te6edf3QueueAdmission
Tff7b72data Tff7b72object T56d364DuplicateId Tb4b4b4: Te6edf3QueueAdmission
Tff7b72data Tff7b72object T56d364ScopeInactive Tb4b4b4: Te6edf3QueueAdmission
Tff7b72data Tff7b72object T56d364TransportUnavailable Tb4b4b4: Te6edf3QueueAdmission
Tb4b4b4}
Tff7b72private Tff7b72val Te6edf3queueResponse Tff7b72= Te6edf3mutableMapOfTff7b72<Tffa657IntTb4b4b4, Te6edf3PendingResponseTff7b72>Tb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3timeoutMutex Tff7b72= Te6edf3MutexTb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3sendAckTimeoutJobs Tff7b72= Te6edf3mutableMapOfTff7b72<Te6edf3PersistedStatusTargetTb4b4b4, Te6edf3JobTff7b72>Tb4b4b4(Tb4b4b4)
Tff7b72override Tff7b72fun Td2a8ffsendToRadioTb4b4b4(Te6edf3pTb4b4b4: Te6edf3ToRadioTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3dispatchToRadioTb4b4b4(Te6edf3pTb4b4b4)Tb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffsendToRadio dropped: no active transport accepted outbound commandTa5d6ff" Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffdispatchToRadioTb4b4b4(Te6edf3pTb4b4b4: Te6edf3ToRadioTb4b4b4)Tb4b4b4: Tffa657Boolean Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffSending to radio Tffd700${Te6edf3pTb4b4b4.Te6edf3toPIIStringTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff" Tb4b4b4}
Tff7b72val Te6edf3dispatched Tff7b72= Te6edf3radioInterfaceServiceTb4b4b4.Te6edf3trySendToRadioTb4b4b4(Te6edf3pTb4b4b4.Te6edf3encodeTb4b4b4(Tb4b4b4)Tb4b4b4)
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3dispatchedTb4b4b4) Tff7b72return Tff7b72false
Te6edf3pTb4b4b4.Te6edf3packetTff7b72?.Te6edf3let Tb4b4b4{ Te6edf3changeStatusTb4b4b4(Tffa657itTb4b4b4, Te6edf3MessageStatusTb4b4b4.Te6edf3ENROUTETb4b4b4) Tb4b4b4}
Tff7b72val Te6edf3packet Tff7b72= Te6edf3pTb4b4b4.Te6edf3packet
Tff7b72if Tb4b4b4(Te6edf3packetTff7b72?.Te6edf3decoded Tff7b72!Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3packetToSave Tff7b72=
Te6edf3MeshLogTb4b4b4(
Te6edf3uuid Tff7b72= Te6edf3UuidTb4b4b4.Te6edf3randomTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3toStringTb4b4b4(Tb4b4b4)Tb4b4b4,
Te6edf3message_type Tff7b72= Ta5d6ff"Ta5d6ffPacketTa5d6ff"Tb4b4b4,
Te6edf3received_date Tff7b72= Te6edf3nowMillisTb4b4b4,
Te6edf3raw_message Tff7b72= Te6edf3packetTb4b4b4.Te6edf3toStringTb4b4b4(Tb4b4b4)Tb4b4b4,
Te6edf3fromNum Tff7b72= Te6edf3MeshLogTb4b4b4.Te6edf3NODE_NUM_LOCALTb4b4b4,
Te6edf3portNum Tff7b72= Te6edf3packetTb4b4b4.Te6edf3decodedTff7b72?.Te6edf3portnumTff7b72?.Te6edf3value Tff7b72?: T79c0ff0Tb4b4b4,
Te6edf3fromRadio Tff7b72= Te6edf3FromRadioTb4b4b4(Te6edf3packet Tff7b72= Te6edf3packetTb4b4b4)Tb4b4b4,
Tb4b4b4)
Te6edf3insertMeshLogTb4b4b4(Te6edf3packetToSaveTb4b4b4)
Tb4b4b4}
Tff7b72return Tff7b72true
Tb4b4b4}
T8b949e/**
* Enqueue [packet] for transmission. Order is preserved for sequential calls from the same coroutine (mutex
* acquisition is uncontested between sequential calls). Transactional sequences that require strict ordering across
* multiple calls — e.g. an `editSettings { … }` begin → writes → commit sequence — MUST be issued from a single
* coroutine; concurrent senders share FIFO only at the per-call grain.
*/
Tff7b72override Tff7b72suspend Tff7b72fun Td2a8ffsendToRadioTb4b4b4(Te6edf3packetTb4b4b4: Te6edf3MeshPacketTb4b4b4)Tb4b4b4: Tffa657Boolean Tff7b72=
Te6edf3enqueueForSendTb4b4b4(Te6edf3packetTb4b4b4, Te6edf3expectedConnectionVersion Tff7b72= Tff7b72nullTb4b4b4)
Tff7b72override Tff7b72suspend Tff7b72fun Td2a8ffsendToRadioForConnectionTb4b4b4(Te6edf3packetTb4b4b4: Te6edf3MeshPacketTb4b4b4, Te6edf3expectedConnectionVersionTb4b4b4: Tffa657LongTb4b4b4)Tb4b4b4: Tffa657Boolean Tff7b72=
Te6edf3enqueueForSendTb4b4b4(Te6edf3packetTb4b4b4, Te6edf3expectedConnectionVersionTb4b4b4)
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffenqueueForSendTb4b4b4(Te6edf3packetTb4b4b4: Te6edf3MeshPacketTb4b4b4, Te6edf3expectedConnectionVersionTb4b4b4: Tffa657Long?Tb4b4b4)Tb4b4b4: Tffa657Boolean Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3packetTb4b4b4.Te6edf3id Tff7b72=Tff7b72= T79c0ff0Tb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffDropping queued packet without an IDTa5d6ff" Tb4b4b4}
Tff7b72return Tff7b72false
Tb4b4b4}
Tff7b72return Tff7b72when Tb4b4b4(Te6edf3enqueuePacketTb4b4b4(Te6edf3packetTb4b4b4, Te6edf3expectedConnectionVersion Tff7b72= Te6edf3expectedConnectionVersionTb4b4b4)Tb4b4b4) Tb4b4b4{
Tff7b72is Te6edf3QueueAdmissionTb4b4b4.Te6edf3Admitted Tff7b72-Tff7b72> Tff7b72true
Te6edf3QueueAdmissionTb4b4b4.Te6edf3TransportUnavailable Tff7b72-Tff7b72> Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffRejecting packet id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff: radio is not connectedTa5d6ff" Tb4b4b4}
Tff7b72false
Tb4b4b4}
Te6edf3QueueAdmissionTb4b4b4.Te6edf3ScopeInactive Tff7b72-Tff7b72> Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffRejecting packet id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff: service scope is no longer activeTa5d6ff" Tb4b4b4}
Tff7b72false
Tb4b4b4}
Te6edf3QueueAdmissionTb4b4b4.Te6edf3DuplicateId Tff7b72-Tff7b72> Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffRejecting packet with reserved id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff" Tb4b4b4}
Tff7b72false
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72suspend Tff7b72fun Td2a8ffsendToRadioAndAwaitResultTb4b4b4(Te6edf3packetTb4b4b4: Te6edf3MeshPacketTb4b4b4)Tb4b4b4: Te6edf3AwaitedSendResult Tff7b72= Tff7b72if Tb4b4b4(Te6edf3packetTb4b4b4.Te6edf3id Tff7b72=Tff7b72= T79c0ff0Tb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffRejecting awaited packet without an IDTa5d6ff" Tb4b4b4}
Te6edf3AwaitedSendResultTb4b4b4(Te6edf3status Tff7b72= Te6edf3AwaitedSendStatusTb4b4b4.Te6edf3REJECTEDTb4b4b4)
Tb4b4b4} Tff7b72else Tb4b4b4{
Tff7b72when Tb4b4b4(Tff7b72val Te6edf3admission Tff7b72= Te6edf3enqueuePacketTb4b4b4(Te6edf3packetTb4b4b4, Te6edf3awaitRouting Tff7b72= Tff7b72trueTb4b4b4)Tb4b4b4) Tb4b4b4{
Tff7b72is Te6edf3QueueAdmissionTb4b4b4.Te6edf3Admitted Tff7b72-Tff7b72> Te6edf3awaitAdmittedPacketTb4b4b4(Te6edf3packetTb4b4b4, Te6edf3admissionTb4b4b4.Te6edf3pendingTb4b4b4)
Te6edf3QueueAdmissionTb4b4b4.Te6edf3TransportUnavailable Tff7b72-Tff7b72> Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffRejecting awaited packet id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff: radio is not connectedTa5d6ff" Tb4b4b4}
Te6edf3AwaitedSendResultTb4b4b4(Te6edf3status Tff7b72= Te6edf3AwaitedSendStatusTb4b4b4.Te6edf3TRANSPORT_STOPPEDTb4b4b4)
Tb4b4b4}
Te6edf3QueueAdmissionTb4b4b4.Te6edf3ScopeInactive Tff7b72-Tff7b72> Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffRejecting awaited packet id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff: service scope is no longer activeTa5d6ff" Tb4b4b4}
Te6edf3AwaitedSendResultTb4b4b4(Te6edf3status Tff7b72= Te6edf3AwaitedSendStatusTb4b4b4.Te6edf3TRANSPORT_STOPPEDTb4b4b4)
Tb4b4b4}
Te6edf3QueueAdmissionTb4b4b4.Te6edf3DuplicateId Tff7b72-Tff7b72> Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffRejecting duplicate awaited packet id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff" Tb4b4b4}
Te6edf3AwaitedSendResultTb4b4b4(Te6edf3status Tff7b72= Te6edf3AwaitedSendStatusTb4b4b4.Te6edf3REJECTEDTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4)
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffawaitAdmittedPacketTb4b4b4(Te6edf3packetTb4b4b4: Te6edf3MeshPacketTb4b4b4, Te6edf3pendingTb4b4b4: Te6edf3PendingResponseTb4b4b4)Tb4b4b4: Te6edf3AwaitedSendResult Tb4b4b4{
Tff7b72val Te6edf3routing Tff7b72= Te6edf3checkNotNullTb4b4b4(Te6edf3pendingTb4b4b4.Te6edf3routingDeferredTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffawaited packet must reserve a routing responseTa5d6ff" Tb4b4b4}
Tff7b72return Tff7b72try Tb4b4b4{
Te6edf3awaitPendingOrServiceStopTb4b4b4(Te6edf3pendingTb4b4b4.Te6edf3dispatchAcceptedTb4b4b4)
Tff7b72val Te6edf3status Tff7b72= Te6edf3awaitPendingOrServiceStopTb4b4b4(Te6edf3routingTb4b4b4)
Te6edf3AwaitedSendResultTb4b4b4(Te6edf3status Tff7b72= Te6edf3statusTb4b4b4, Te6edf3departureEpochAtDispatch Tff7b72= Te6edf3pendingTb4b4b4.Te6edf3departureEpochAtDispatchTb4b4b4)
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3CancellationExceptionTb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3e T8b949e// Caller cancellation does not release queue or response ownership.
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3eTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffsendToRadioAndAwait packet id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff failedTa5d6ff" Tb4b4b4}
Tff7b72val Te6edf3failureStatus Tff7b72= Te6edf3pendingTb4b4b4.Te6edf3sendFailureStatusTb4b4b4(Tb4b4b4)
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3pendingTb4b4b4.Te6edf3completeRoutingTb4b4b4(Te6edf3failureStatusTb4b4b4)
Te6edf3removePendingIfTerminalLockedTb4b4b4(Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4, Te6edf3pendingTb4b4b4)
Tb4b4b4}
Te6edf3AwaitedSendResultTb4b4b4(
Te6edf3status Tff7b72= Te6edf3pendingTb4b4b4.Te6edf3terminalStatus Tff7b72?: Te6edf3failureStatusTb4b4b4,
Te6edf3departureEpochAtDispatch Tff7b72= Te6edf3pendingTb4b4b4.Te6edf3departureEpochAtDispatchTb4b4b4,
Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Reserves [packet]'s non-zero ID and queues it as one atomic admission. The reservation spans queued work, the
* firmware queue result, and (for strict callers) the later routing acknowledgement.
*/
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffenqueuePacketTb4b4b4(
Te6edf3packetTb4b4b4: Te6edf3MeshPacketTb4b4b4,
Te6edf3awaitRoutingTb4b4b4: Tffa657Boolean Tff7b72= Tff7b72falseTb4b4b4,
Te6edf3expectedConnectionVersionTb4b4b4: Tffa657Long? Tff7b72= Tff7b72nullTb4b4b4,
Tb4b4b4)Tb4b4b4: Te6edf3QueueAdmission Tff7b72= Te6edf3queueMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tf0883eresponseLock@Tb4b4b4{
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3scopeTb4b4b4.Te6edf3isActiveTb4b4b4) Tff7b72returnTf0883e@responseLock Te6edf3QueueAdmissionTb4b4b4.Te6edf3ScopeInactive
Tff7b72if Tb4b4b4(Te6edf3expectedConnectionVersion Tff7b72=Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3connectionStateProviderTb4b4b4.Te6edf3connectionStateTb4b4b4.Te6edf3value Tff7b72!Tff7b72= Te6edf3ConnectionStateTb4b4b4.Te6edf3ConnectedTb4b4b4) Tb4b4b4{
Tff7b72returnTf0883e@responseLock Te6edf3QueueAdmissionTb4b4b4.Te6edf3TransportUnavailable
Tb4b4b4}
Tb4b4b4} Tff7b72else Tb4b4b4{
Tff7b72val Te6edf3lifecycle Tff7b72= Te6edf3connectionStateProviderTb4b4b4.Te6edf3connectionLifecycleTb4b4b4.Te6edf3value
Tff7b72if Tb4b4b4(
Te6edf3lifecycleTb4b4b4.Te6edf3state Tff7b72!is Te6edf3ConnectionStateTb4b4b4.Te6edf3Connected Tff7b72|Tff7b72| Te6edf3lifecycleTb4b4b4.Te6edf3version Tff7b72!Tff7b72= Te6edf3expectedConnectionVersion
Tb4b4b4) Tb4b4b4{
Tff7b72returnTf0883e@responseLock Te6edf3QueueAdmissionTb4b4b4.Te6edf3TransportUnavailable
Tb4b4b4}
Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3queueResponseTb4b4b4.Te6edf3containsKeyTb4b4b4(Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4)Tb4b4b4) Tff7b72returnTf0883e@responseLock Te6edf3QueueAdmissionTb4b4b4.Te6edf3DuplicateId
Tff7b72val Te6edf3pending Tff7b72=
Te6edf3PendingResponseTb4b4b4(
Te6edf3packet Tff7b72= Te6edf3packetTb4b4b4,
Te6edf3persistence Tff7b72= Te6edf3packetTb4b4b4.Te6edf3statusPersistenceTb4b4b4(Tb4b4b4)Tb4b4b4,
Te6edf3awaitRouting Tff7b72= Te6edf3awaitRoutingTb4b4b4,
Tb4b4b4)
Te6edf3queueResponseTff7b72[Te6edf3packetTb4b4b4.Te6edf3idTff7b72] Tff7b72= Te6edf3pending
Te6edf3queueStopped Tff7b72= Tff7b72false T8b949e// Allow queue to resume after a disconnect/reconnect cycle.
Te6edf3queuedPacketsTb4b4b4.Te6edf3addTb4b4b4(
Te6edf3QueuedPacketTb4b4b4(
Te6edf3packet Tff7b72= Te6edf3packetTb4b4b4,
Te6edf3pending Tff7b72= Te6edf3pendingTb4b4b4,
Te6edf3expectedConnectionVersion Tff7b72= Te6edf3expectedConnectionVersionTb4b4b4,
Tb4b4b4)Tb4b4b4,
Tb4b4b4)
Te6edf3startPacketQueueLockedTb4b4b4(Tb4b4b4)
Te6edf3QueueAdmissionTb4b4b4.Te6edf3AdmittedTb4b4b4(Te6edf3pendingTb4b4b4)
Tb4b4b4}
Tb4b4b4}
T8b949e/** Waits for one pending stage while converting permanent service shutdown into a synchronous queue drain. */
Tff7b72private Tff7b72suspend Tff7b72fun Tff7b72<Te6edf3TTff7b72> Td2a8ffawaitPendingOrServiceStopTb4b4b4(Te6edf3deferredTb4b4b4: Te6edf3DeferredTff7b72<Te6edf3TTff7b72>Tb4b4b4)Tb4b4b4: Te6edf3T Tb4b4b4{
Tff7b72val Te6edf3scopeJob Tff7b72= Te6edf3scopeTb4b4b4.Te6edf3coroutineContextTff7b72[Te6edf3JobTff7b72] Tff7b72?: Tff7b72return Te6edf3deferredTb4b4b4.Te6edf3awaitTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3serviceStopped Tff7b72= Te6edf3select Tb4b4b4{
Te6edf3deferredTb4b4b4.Te6edf3onAwait Tb4b4b4{ Tff7b72false Tb4b4b4}
Te6edf3scopeJobTb4b4b4.Te6edf3onJoin Tb4b4b4{ Tff7b72true Tb4b4b4}
Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3serviceStoppedTb4b4b4) Tb4b4b4{
Te6edf3withContextTb4b4b4(Te6edf3NonCancellableTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3failedPacketIds Tff7b72=
Te6edf3queueMutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Tff7b72if Tb4b4b4(Tff7b72!Te6edf3deferredTb4b4b4.Te6edf3isCompletedTb4b4b4) Te6edf3stopAndDrainPacketQueueLockedTb4b4b4(Tb4b4b4) Tff7b72else Te6edf3emptyListTb4b4b4(Tb4b4b4) Tb4b4b4}
Te6edf3changeStatusesNowTb4b4b4(Te6edf3failedPacketIdsTb4b4b4, Te6edf3MessageStatusTb4b4b4.Te6edf3ERRORTb4b4b4, Te6edf3PERSISTED_STATUS_SHUTDOWN_WAITTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tff7b72return Te6edf3deferredTb4b4b4.Te6edf3awaitTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffstopPacketQueueTb4b4b4(Tb4b4b4) Tb4b4b4{
T8b949e// Run async so callers (non-suspend) don't block, but all mutations are
T8b949e// serialized under the same mutexes used by the queue processor and senders.
Te6edf3scopeTb4b4b4.Te6edf3handledLaunch Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{ Ta5d6ff"Ta5d6ffStopping packet queueJobTa5d6ff" Tb4b4b4}
Te6edf3withContextTb4b4b4(Te6edf3NonCancellableTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3failedPacketIds Tff7b72=
Te6edf3queueMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3queueStopped Tff7b72= Tff7b72true
Te6edf3queueJobTff7b72?.Te6edf3cancelTb4b4b4(Tb4b4b4)
Te6edf3queueJob Tff7b72= Tff7b72null
Te6edf3queueGenerationTff7b72+Tff7b72+
Te6edf3queuedPacketsTb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)
Te6edf3completePendingResponsesTb4b4b4(Te6edf3AwaitedSendStatusTb4b4b4.Te6edf3TRANSPORT_STOPPEDTb4b4b4)
Tb4b4b4}
Te6edf3changeStatusesNowTb4b4b4(Te6edf3failedPacketIdsTb4b4b4, Te6edf3MessageStatusTb4b4b4.Te6edf3ERRORTb4b4b4, Te6edf3PERSISTED_STATUS_SHUTDOWN_WAITTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffhandleQueueStatusTb4b4b4(Te6edf3queueStatusTb4b4b4: Te6edf3QueueStatusTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ff[queueStatus] Tffd700${Te6edf3queueStatusTb4b4b4.Te6edf3toOneLineStringTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff" Tb4b4b4}
Tff7b72val Tb4b4b4(Te6edf3successTb4b4b4, Te6edf3isFullTb4b4b4, Te6edf3requestIdTb4b4b4) Tff7b72=
Te6edf3withTb4b4b4(Te6edf3queueStatusTb4b4b4) Tb4b4b4{ Te6edf3TripleTb4b4b4(Te6edf3res Tff7b72=Tff7b72= T79c0ff0 Tff7b72|Tff7b72| Te6edf3res Tff7b72=Tff7b72= Te6edf3ERRNO_SHOULD_RELEASETb4b4b4, Te6edf3free Tff7b72=Tff7b72= T79c0ff0Tb4b4b4, Te6edf3mesh_packet_idTb4b4b4) Tb4b4b4}
T8b949e// Only plain res==0 "accepted, now full" is an advisory echo. ERRNO_SHOULD_RELEASE is synchronous local
T8b949e// delivery and therefore completes both the firmware-queue stage and any strict routing waiter.
Tff7b72if Tb4b4b4(Te6edf3queueStatusTb4b4b4.Te6edf3res Tff7b72=Tff7b72= T79c0ff0 Tff7b72&Tff7b72& Te6edf3isFull Tff7b72&Tff7b72& Te6edf3requestId Tff7b72=Tff7b72= T79c0ff0Tb4b4b4) Tff7b72return
Te6edf3scopeTb4b4b4.Te6edf3handledLaunch Tb4b4b4{
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72val Te6edf3entry Tff7b72=
Tff7b72if Tb4b4b4(Te6edf3requestId Tff7b72!Tff7b72= T79c0ff0Tb4b4b4) Tb4b4b4{
Te6edf3queueResponseTff7b72[Te6edf3requestIdTff7b72]Tff7b72?.Te6edf3let Tb4b4b4{ Te6edf3requestId Te6edf3to Tffa657it Tb4b4b4}
Tb4b4b4} Tff7b72else Tb4b4b4{
T8b949e// Firmware can omit mesh_packet_id for the active queue entry. Strict routing waiters whose
T8b949e// queue stage already completed remain in the map, so correlate only an unfinished queue stage.
Te6edf3queueResponseTb4b4b4.Te6edf3entries
Tb4b4b4.Te6edf3firstOrNull Tb4b4b4{ Tb4b4b4(Te6edf3_Tb4b4b4, Te6edf3pendingTb4b4b4) Tff7b72-Tff7b72> Te6edf3pendingTb4b4b4.Te6edf3wasDispatched Tff7b72&Tff7b72& Tff7b72!Te6edf3pendingTb4b4b4.Te6edf3queueDeferredTb4b4b4.Te6edf3isCompleted Tb4b4b4}
Tff7b72?.Te6edf3let Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3key Te6edf3to Tffa657itTb4b4b4.Te6edf3value Tb4b4b4}
Tb4b4b4}
Tff7b72val Tb4b4b4(Te6edf3packetIdTb4b4b4, Te6edf3pendingTb4b4b4) Tff7b72= Te6edf3entry Tff7b72?: Tff7b72returnTf0883e@withLock
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3pendingTb4b4b4.Te6edf3wasDispatchedTb4b4b4) Tff7b72returnTf0883e@withLock
Tff7b72val Te6edf3status Tff7b72= Te6edf3successTb4b4b4.Te6edf3toAwaitedSendStatusTb4b4b4(Tb4b4b4)
Te6edf3pendingTb4b4b4.Te6edf3completeQueueTb4b4b4(Te6edf3statusTb4b4b4, Te6edf3routingTerminal Tff7b72= Tff7b72!Te6edf3success Tff7b72|Tff7b72| Te6edf3queueStatusTb4b4b4.Te6edf3res Tff7b72=Tff7b72= Te6edf3ERRNO_SHOULD_RELEASETb4b4b4)
Te6edf3removePendingIfTerminalLockedTb4b4b4(Te6edf3packetIdTb4b4b4, Te6edf3pendingTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72suspend Tff7b72fun Td2a8ffcompleteDispatchedResponseTb4b4b4(Te6edf3dataRequestIdTb4b4b4: Tffa657IntTb4b4b4, Te6edf3completeTb4b4b4: Tffa657BooleanTb4b4b4) Tb4b4b4{
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72val Te6edf3pending Tff7b72= Te6edf3queueResponseTff7b72[Te6edf3dataRequestIdTff7b72]Tff7b72?.Te6edf3takeIf Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3wasDispatched Tb4b4b4} Tff7b72?: Tff7b72returnTf0883e@withLock
Te6edf3pendingTb4b4b4.Te6edf3completeRoutingTb4b4b4(Te6edf3completeTb4b4b4.Te6edf3toAwaitedSendStatusTb4b4b4(Tb4b4b4)Tb4b4b4)
Te6edf3removePendingIfTerminalLockedTb4b4b4(Te6edf3dataRequestIdTb4b4b4, Te6edf3pendingTb4b4b4)
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Starts the packet queue processor. Must be called while holding [queueMutex] to ensure the check-then-start is
* atomic — preventing two concurrent callers from launching duplicate processors.
*/
Tff7b72private Tff7b72fun Td2a8ffstartPacketQueueLockedTb4b4b4(Tb4b4b4) Tb4b4b4{
Te6edf3checkTb4b4b4(Te6edf3queueMutexTb4b4b4.Te6edf3isLockedTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffPacket queue workers must start while queueMutex is heldTa5d6ff" Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3queueStopped Tff7b72|Tff7b72| Te6edf3queueJobTff7b72?.Te6edf3isActive Tff7b72=Tff7b72= Tff7b72trueTb4b4b4) Tff7b72return
Tff7b72val Te6edf3generation Tff7b72= Tff7b72+Tff7b72+Te6edf3queueGeneration
T8b949e// Install cleanup before admission releases queueMutex so cancellation cannot bypass the worker's finally.
T8b949e// UNDISPATCHED enters immediately, then suspends on queueMutex until enqueuePacket releases the admission
T8b949e// lock; this closes the cancellation gap without running queue mutation inside the caller's critical section.
Te6edf3queueJob Tff7b72=
Te6edf3scopeTb4b4b4.Te6edf3handledLaunchTb4b4b4(Te6edf3start Tff7b72= Te6edf3CoroutineStartTb4b4b4.Te6edf3UNDISPATCHEDTb4b4b4) Tb4b4b4{
Tff7b72try Tb4b4b4{
Tff7b72while Tb4b4b4(Te6edf3connectionStateProviderTb4b4b4.Te6edf3connectionStateTb4b4b4.Te6edf3value Tff7b72=Tff7b72= Te6edf3ConnectionStateTb4b4b4.Te6edf3ConnectedTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3queuedPacket Tff7b72= Te6edf3queueMutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Te6edf3queuedPacketsTb4b4b4.Te6edf3removeFirstOrNullTb4b4b4(Tb4b4b4) Tb4b4b4} Tff7b72?: Tff7b72break
Te6edf3processQueuedPacketTb4b4b4(Te6edf3queuedPacketTb4b4b4)
Tb4b4b4}
Tb4b4b4} Tff7b72finally Tb4b4b4{
Te6edf3finishPacketQueueGenerationTb4b4b4(Te6edf3generationTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4)
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffprocessQueuedPacketTb4b4b4(Te6edf3queuedPacketTb4b4b4: Te6edf3QueuedPacketTb4b4b4) Tb4b4b4{
Tff7b72val Tb4b4b4(Te6edf3packetTb4b4b4, Te6edf3pendingTb4b4b4, Te6edf3expectedConnectionVersionTb4b4b4) Tff7b72= Te6edf3queuedPacket
Tff7b72try Tb4b4b4{
Tff7b72val Te6edf3response Tff7b72= Te6edf3sendPacketTb4b4b4(Te6edf3packetTb4b4b4, Te6edf3pendingTb4b4b4, Te6edf3expectedConnectionVersionTb4b4b4)
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffqueueJob packet id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff waiting for QueueStatusTa5d6ff" Tb4b4b4}
Tff7b72val Te6edf3status Tff7b72= Te6edf3withTimeoutTb4b4b4(Te6edf3RESPONSE_TIMEOUTTb4b4b4) Tb4b4b4{ Te6edf3responseTb4b4b4.Te6edf3awaitTb4b4b4(Tb4b4b4) Tb4b4b4}
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffqueueJob packet id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff queue status Tffd700$Te6edf3statusTa5d6ff" Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3status Tff7b72!Tff7b72= Te6edf3AwaitedSendStatusTb4b4b4.Te6edf3ACCEPTEDTb4b4b4) Te6edf3changeStatusTb4b4b4(Te6edf3packetTb4b4b4, Te6edf3MessageStatusTb4b4b4.Te6edf3ERRORTb4b4b4)
Te6edf3removePendingIfTerminalTb4b4b4(Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4, Te6edf3pendingTb4b4b4)
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3_Tb4b4b4: Te6edf3TimeoutCancellationExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffqueueJob packet id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff queue response timeoutTa5d6ff" Tb4b4b4}
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
T8b949e// QueueStatus is a bounded, lossy phone-side signal in firmware. A timeout after transport dispatch is
T8b949e// therefore inconclusive: release the local queue stage, but keep any strict routing waiter alive.
Te6edf3pendingTb4b4b4.Te6edf3completeQueueTb4b4b4(Te6edf3AwaitedSendStatusTb4b4b4.Te6edf3TIMED_OUTTb4b4b4, Te6edf3routingTerminal Tff7b72= Tff7b72falseTb4b4b4)
Te6edf3removePendingIfTerminalLockedTb4b4b4(Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4, Te6edf3pendingTb4b4b4)
Tb4b4b4}
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3CancellationExceptionTb4b4b4) Tb4b4b4{
Te6edf3withContextTb4b4b4(Te6edf3NonCancellableTb4b4b4) Tb4b4b4{
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3pendingTb4b4b4.Te6edf3completeAllTb4b4b4(Te6edf3AwaitedSendStatusTb4b4b4.Te6edf3TRANSPORT_STOPPEDTb4b4b4)
Te6edf3removePendingIfTerminalLockedTb4b4b4(Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4, Te6edf3pendingTb4b4b4)
Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3pendingTb4b4b4.Te6edf3terminalStatus Tff7b72!Tff7b72= Te6edf3AwaitedSendStatusTb4b4b4.Te6edf3ACCEPTEDTb4b4b4) Tb4b4b4{
Te6edf3changeStatusNowTb4b4b4(Te6edf3packetTb4b4b4, Te6edf3MessageStatusTb4b4b4.Te6edf3ERRORTb4b4b4, Te6edf3PERSISTED_STATUS_SHUTDOWN_WAITTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tff7b72throw Te6edf3e T8b949e// Preserve structured concurrency cancellation propagation.
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3eTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffqueueJob packet id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff failedTa5d6ff" Tb4b4b4}
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3pendingTb4b4b4.Te6edf3completeAllTb4b4b4(Te6edf3pendingTb4b4b4.Te6edf3sendFailureStatusTb4b4b4(Tb4b4b4)Tb4b4b4)
Te6edf3removePendingIfTerminalLockedTb4b4b4(Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4, Te6edf3pendingTb4b4b4)
Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3pendingTb4b4b4.Te6edf3terminalStatus Tff7b72!Tff7b72= Te6edf3AwaitedSendStatusTb4b4b4.Te6edf3ACCEPTEDTb4b4b4) Te6edf3changeStatusTb4b4b4(Te6edf3packetTb4b4b4, Te6edf3MessageStatusTb4b4b4.Te6edf3ERRORTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8fffinishPacketQueueGenerationTb4b4b4(Te6edf3generationTb4b4b4: Tffa657LongTb4b4b4) Tff7b72= Te6edf3withContextTb4b4b4(Te6edf3NonCancellableTb4b4b4) Tb4b4b4{
T8b949e// Keep completion, replacement, and disconnect draining atomic with new admissions. queueGeneration
T8b949e// advances only under queueMutex: stopPacketQueue() drains every pending response, while
T8b949e// startPacketQueueLocked() advances it only after the previous queueJob is inactive. Therefore a stale
T8b949e// worker has already lost ownership to a path that drained or replaced it and must not clear its successor.
Tff7b72val Te6edf3failedPacketIds Tff7b72=
Te6edf3queueMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3generation Tff7b72!Tff7b72= Te6edf3queueGenerationTb4b4b4) Tff7b72returnTf0883e@withLock Te6edf3emptyListTb4b4b4(Tb4b4b4)
Te6edf3queueJob Tff7b72= Tff7b72null
Tff7b72when Tb4b4b4{
Te6edf3queueStopped Tff7b72|Tff7b72| Tff7b72!Te6edf3scopeTb4b4b4.Te6edf3isActive Tff7b72-Tff7b72> Te6edf3stopAndDrainPacketQueueLockedTb4b4b4(Tb4b4b4)
Te6edf3connectionStateProviderTb4b4b4.Te6edf3connectionStateTb4b4b4.Te6edf3value Tff7b72!Tff7b72= Te6edf3ConnectionStateTb4b4b4.Te6edf3Connected Tff7b72-Tff7b72>
Te6edf3stopAndDrainPacketQueueLockedTb4b4b4(Tb4b4b4)
Te6edf3queuedPacketsTb4b4b4.Te6edf3isNotEmptyTb4b4b4(Tb4b4b4) Tff7b72-Tff7b72> Tb4b4b4{
Te6edf3startPacketQueueLockedTb4b4b4(Tb4b4b4)
Te6edf3emptyListTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tff7b72else Tff7b72-Tff7b72> Te6edf3emptyListTb4b4b4(Tb4b4b4) T8b949e// Strict routing waiters may remain after their QueueStatus completed.
Tb4b4b4}
Tb4b4b4}
Te6edf3changeStatusesNowTb4b4b4(Te6edf3failedPacketIdsTb4b4b4, Te6edf3MessageStatusTb4b4b4.Te6edf3ERRORTb4b4b4, Te6edf3PERSISTED_STATUS_SHUTDOWN_WAITTb4b4b4)
Tb4b4b4}
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffstopAndDrainPacketQueueLockedTb4b4b4(Tb4b4b4)Tb4b4b4: Te6edf3ListTff7b72<Te6edf3PacketStatusTargetTff7b72> Tb4b4b4{
Te6edf3queueStopped Tff7b72= Tff7b72true
Te6edf3queuedPacketsTb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)
Tff7b72return Te6edf3completePendingResponsesTb4b4b4(Te6edf3AwaitedSendStatusTb4b4b4.Te6edf3TRANSPORT_STOPPEDTb4b4b4)
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffrearmSendAckTimeoutsTb4b4b4(Tb4b4b4) Tb4b4b4{
Te6edf3scopeTb4b4b4.Te6edf3handledLaunch Tb4b4b4{
Te6edf3packetRepositoryTb4b4b4.Te6edf3valueTb4b4b4.Te6edf3getEnroutePacketsTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3forEach Tb4b4b4{ Te6edf3persisted Tff7b72-Tff7b72>
Tff7b72val Te6edf3remaining Tff7b72= Te6edf3persistedTb4b4b4.Te6edf3packetTb4b4b4.Te6edf3time Tff7b72+ Te6edf3SEND_ACK_TIMEOUTTb4b4b4.Te6edf3inWholeMilliseconds Tff7b72- Te6edf3nowMillis
Te6edf3scheduleSendAckTimeoutTb4b4b4(
Te6edf3PersistedStatusTargetTb4b4b4.Te6edf3DataPacketTb4b4b4(Te6edf3persistedTb4b4b4.Te6edf3idTb4b4b4)Tb4b4b4,
Te6edf3remainingTb4b4b4.Te6edf3millisecondsTb4b4b4.Te6edf3coerceAtLeastTb4b4b4(Te6edf3REARM_GRACETb4b4b4)Tb4b4b4,
Tb4b4b4)
Tb4b4b4}
Te6edf3packetRepositoryTb4b4b4.Te6edf3valueTb4b4b4.Te6edf3getEnrouteReactionsTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3forEach Tb4b4b4{ Te6edf3persisted Tff7b72-Tff7b72>
Tff7b72val Te6edf3remaining Tff7b72= Te6edf3persistedTb4b4b4.Te6edf3reactionTb4b4b4.Te6edf3timestamp Tff7b72+ Te6edf3SEND_ACK_TIMEOUTTb4b4b4.Te6edf3inWholeMilliseconds Tff7b72- Te6edf3nowMillis
Te6edf3scheduleSendAckTimeoutTb4b4b4(
Te6edf3PersistedStatusTargetTb4b4b4.Te6edf3ReactionTb4b4b4(Te6edf3persistedTb4b4b4.Te6edf3idTb4b4b4)Tb4b4b4,
Te6edf3remainingTb4b4b4.Te6edf3millisecondsTb4b4b4.Te6edf3coerceAtLeastTb4b4b4(Te6edf3REARM_GRACETb4b4b4)Tb4b4b4,
Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* A send whose routing ACK/NAK never reaches the app (typically because it was disconnected when the radio's
* response arrived) would stay ENROUTE — "Sending…" — forever. Stamp it as a retryable timeout instead; a late ACK
* still upgrades it via handleAckNak.
*
* Packet and reaction timers are keyed by stable persisted-row identity. Re-arming supersedes the pending timer for
* the same target, and timers deliberately survive disconnects.
*/
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffscheduleSendAckTimeoutTb4b4b4(Te6edf3targetTb4b4b4: Te6edf3PersistedStatusTargetTb4b4b4, Te6edf3delayForTb4b4b4: Te6edf3Duration Tff7b72= Te6edf3SEND_ACK_TIMEOUTTb4b4b4) Tb4b4b4{
Te6edf3timeoutMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3sendAckTimeoutJobsTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3targetTb4b4b4)Tff7b72?.Te6edf3cancelTb4b4b4(Tb4b4b4)
Te6edf3sendAckTimeoutJobsTb4b4b4.Te6edf3valuesTb4b4b4.Te6edf3removeAll Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3isCompleted Tb4b4b4}
Te6edf3sendAckTimeoutJobsTff7b72[Te6edf3targetTff7b72] Tff7b72=
Te6edf3scopeTb4b4b4.Te6edf3handledLaunch Tb4b4b4{
Te6edf3delayTb4b4b4(Te6edf3delayForTb4b4b4)
T8b949e// Conditional in the DAO transaction: an ACK/NAK landing while this timer waited must win.
Tff7b72when Tb4b4b4(Te6edf3targetTb4b4b4) Tb4b4b4{
Tff7b72is Te6edf3PersistedStatusTargetTb4b4b4.Te6edf3DataPacket Tff7b72-Tff7b72>
Te6edf3packetRepositoryTb4b4b4.Te6edf3valueTb4b4b4.Te6edf3timeOutEnroutePacketTb4b4b4(Te6edf3targetTb4b4b4.Te6edf3idTb4b4b4, Te6edf3RoutingTb4b4b4.Te6edf3ErrorTb4b4b4.Te6edf3TIMEOUTTb4b4b4.Te6edf3valueTb4b4b4)
Tff7b72is Te6edf3PersistedStatusTargetTb4b4b4.Te6edf3Reaction Tff7b72-Tff7b72>
Te6edf3packetRepositoryTb4b4b4.Te6edf3valueTb4b4b4.Te6edf3timeOutEnrouteReactionTb4b4b4(Te6edf3targetTb4b4b4.Te6edf3idTb4b4b4, Te6edf3RoutingTb4b4b4.Te6edf3ErrorTb4b4b4.Te6edf3TIMEOUTTb4b4b4.Te6edf3valueTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffchangeStatusTb4b4b4(Te6edf3packetTb4b4b4: Te6edf3MeshPacketTb4b4b4, Te6edf3statusTb4b4b4: Te6edf3MessageStatusTb4b4b4) Tff7b72=
Te6edf3scopeTb4b4b4.Te6edf3handledLaunch Tb4b4b4{ Te6edf3changeStatusNowTb4b4b4(Te6edf3packetTb4b4b4, Te6edf3statusTb4b4b4) Tb4b4b4}
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffchangeStatusNowTb4b4b4(
Te6edf3packetTb4b4b4: Te6edf3MeshPacketTb4b4b4,
Te6edf3statusTb4b4b4: Te6edf3MessageStatusTb4b4b4,
Te6edf3persistenceWaitTb4b4b4: Te6edf3Duration Tff7b72= Te6edf3PERSISTED_STATUS_WAITTb4b4b4,
Tb4b4b4) Tb4b4b4{
Te6edf3changeStatusesNowTb4b4b4(
Te6edf3targets Tff7b72= Te6edf3listOfTb4b4b4(Te6edf3PacketStatusTargetTb4b4b4(Te6edf3packetTb4b4b4, Te6edf3packetTb4b4b4.Te6edf3statusPersistenceTb4b4b4(Tb4b4b4)Tb4b4b4)Tb4b4b4)Tb4b4b4,
Te6edf3status Tff7b72= Te6edf3statusTb4b4b4,
Te6edf3persistenceWait Tff7b72= Te6edf3persistenceWaitTb4b4b4,
Tb4b4b4)
Tb4b4b4}
T8b949e/**
* Applies [status] to durable rows, sharing one brief wait across app payloads whose inserts may still be in
* flight. Data packets resolve by the full outgoing identity from #6624 before their exact persisted row is
* mutated, so sender-scoped packet-ID collisions and terminal ACK/NAK states both fail closed.
*/
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffchangeStatusesNowTb4b4b4(
Te6edf3targetsTb4b4b4: Te6edf3CollectionTff7b72<Te6edf3PacketStatusTargetTff7b72>Tb4b4b4,
Te6edf3statusTb4b4b4: Te6edf3MessageStatusTb4b4b4,
Te6edf3persistenceWaitTb4b4b4: Te6edf3Duration Tff7b72= Te6edf3PERSISTED_STATUS_WAITTb4b4b4,
Tb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3waiting Tff7b72=
Te6edf3targets
Tb4b4b4.Te6edf3filter Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3packetTb4b4b4.Te6edf3id Tff7b72!Tff7b72= T79c0ff0 Tff7b72&Tff7b72& Tffa657itTb4b4b4.Te6edf3persistence Tff7b72!Tff7b72= Tff7b72null Tb4b4b4}
Tb4b4b4.Te6edf3distinctBy Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3packetTb4b4b4.Te6edf3id Tb4b4b4}
Tb4b4b4.Te6edf3associateByToTb4b4b4(Te6edf3mutableMapOfTb4b4b4(Tb4b4b4)Tb4b4b4) Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3packetTb4b4b4.Te6edf3id Tb4b4b4}
Te6edf3withTimeoutOrNullTb4b4b4(Te6edf3persistenceWaitTb4b4b4) Tb4b4b4{
Tff7b72while Tb4b4b4(Te6edf3waitingTb4b4b4.Te6edf3isNotEmptyTb4b4b4(Tb4b4b4)Tb4b4b4) Tb4b4b4{
Te6edf3waitingTb4b4b4.Te6edf3valuesTb4b4b4.Te6edf3toListTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3forEach Tb4b4b4{ Te6edf3target Tff7b72-Tff7b72>
Tff7b72if Tb4b4b4(Te6edf3applyQueueStatusTb4b4b4(Te6edf3targetTb4b4b4, Te6edf3statusTb4b4b4)Tb4b4b4) Te6edf3waitingTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3targetTb4b4b4.Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4)
Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3waitingTb4b4b4.Te6edf3isNotEmptyTb4b4b4(Tb4b4b4)Tb4b4b4) Te6edf3delayTb4b4b4(Te6edf3PERSISTED_STATUS_RETRY_DELAYTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Te6edf3waitingTb4b4b4.Te6edf3keysTb4b4b4.Te6edf3forEach Tb4b4b4{ Te6edf3packetId Tff7b72-Tff7b72> Te6edf3logMissingStatusRowTb4b4b4(Te6edf3packetIdTb4b4b4, Te6edf3statusTb4b4b4) Tb4b4b4}
Tb4b4b4}
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffapplyQueueStatusTb4b4b4(Te6edf3targetTb4b4b4: Te6edf3PacketStatusTargetTb4b4b4, Te6edf3statusTb4b4b4: Te6edf3MessageStatusTb4b4b4)Tb4b4b4: Tffa657Boolean Tff7b72=
Tff7b72when Tb4b4b4(Te6edf3targetTb4b4b4.Te6edf3persistenceTb4b4b4) Tb4b4b4{
Te6edf3StatusPersistenceTb4b4b4.Te6edf3DATA_PACKET Tff7b72-Tff7b72> Tb4b4b4{
Tff7b72val Te6edf3persisted Tff7b72= Te6edf3packetRepositoryTb4b4b4.Te6edf3valueTb4b4b4.Te6edf3applyOutgoingQueueStatusTb4b4b4(Te6edf3targetTb4b4b4.Te6edf3packetTb4b4b4, Te6edf3statusTb4b4b4) Tff7b72?: Tff7b72return Tff7b72false
Tff7b72val Te6edf3current Tff7b72= Te6edf3persistedTb4b4b4.Te6edf3packetTb4b4b4.Te6edf3status
Tff7b72val Te6edf3shouldPersist Tff7b72=
Te6edf3status Tff7b72=Tff7b72= Te6edf3MessageStatusTb4b4b4.Te6edf3ENROUTE Tff7b72&Tff7b72&
Tb4b4b4(Te6edf3current Tff7b72=Tff7b72= Te6edf3MessageStatusTb4b4b4.Te6edf3ENROUTE Tff7b72|Tff7b72| Te6edf3shouldApplyOutgoingQueueStatusTb4b4b4(Te6edf3currentTb4b4b4, Te6edf3statusTb4b4b4)Tb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3shouldPersistTb4b4b4) Tb4b4b4{
Te6edf3scheduleSendAckTimeoutTb4b4b4(Te6edf3PersistedStatusTargetTb4b4b4.Te6edf3DataPacketTb4b4b4(Te6edf3persistedTb4b4b4.Te6edf3idTb4b4b4)Tb4b4b4)
Tb4b4b4}
Tff7b72true
Tb4b4b4}
Te6edf3StatusPersistenceTb4b4b4.Te6edf3REACTION Tff7b72-Tff7b72> Tb4b4b4{
Tff7b72val Te6edf3persisted Tff7b72=
Te6edf3packetRepositoryTb4b4b4.Te6edf3valueTb4b4b4.Te6edf3applyOutgoingReactionQueueStatusTb4b4b4(Te6edf3targetTb4b4b4.Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4, Te6edf3statusTb4b4b4) Tff7b72?: Tff7b72return Tff7b72false
Tff7b72val Te6edf3current Tff7b72= Te6edf3persistedTb4b4b4.Te6edf3reactionTb4b4b4.Te6edf3status
Tff7b72if Tb4b4b4(
Te6edf3status Tff7b72=Tff7b72= Te6edf3MessageStatusTb4b4b4.Te6edf3ENROUTE Tff7b72&Tff7b72&
Tb4b4b4(Te6edf3current Tff7b72=Tff7b72= Te6edf3MessageStatusTb4b4b4.Te6edf3ENROUTE Tff7b72|Tff7b72| Te6edf3shouldApplyOutgoingQueueStatusTb4b4b4(Te6edf3currentTb4b4b4, Te6edf3statusTb4b4b4)Tb4b4b4)
Tb4b4b4) Tb4b4b4{
Te6edf3scheduleSendAckTimeoutTb4b4b4(Te6edf3PersistedStatusTargetTb4b4b4.Te6edf3ReactionTb4b4b4(Te6edf3persistedTb4b4b4.Te6edf3idTb4b4b4)Tb4b4b4)
Tb4b4b4}
Tff7b72true
Tb4b4b4}
Tff7b72null Tff7b72-Tff7b72> Tff7b72true
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8fflogMissingStatusRowTb4b4b4(Te6edf3packetIdTb4b4b4: Tffa657IntTb4b4b4, Te6edf3statusTb4b4b4: Te6edf3MessageStatusTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffSkipping Tffd700$Te6edf3statusTa5d6ff for mesh packet id=Tffd700${Te6edf3packetIdTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff: no unambiguous persisted rowTa5d6ff" Tb4b4b4}
Tb4b4b4}
Tff7b72private Tff7b72fun Te6edf3MeshPacketTb4b4b4.Td2a8ffstatusPersistenceTb4b4b4(Tb4b4b4)Tb4b4b4: Te6edf3StatusPersistence? Tb4b4b4{
Tff7b72val Te6edf3data Tff7b72= Te6edf3decoded Tff7b72?: Tff7b72return Tff7b72null
Tff7b72return Tff7b72when Tb4b4b4{
Te6edf3dataTb4b4b4.Te6edf3isReactionTb4b4b4(Tb4b4b4) Tff7b72-Tff7b72> Te6edf3StatusPersistenceTb4b4b4.Te6edf3REACTION
Te6edf3dataTb4b4b4.Te6edf3portnumTb4b4b4.Te6edf3value Tff7b72in Te6edf3PERSISTED_DATA_PORT_NUMBERS Tff7b72-Tff7b72> Te6edf3StatusPersistenceTb4b4b4.Te6edf3DATA_PACKET
Tff7b72else Tff7b72-Tff7b72> Tff7b72null
Tb4b4b4}
Tb4b4b4}
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4)
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffsendPacketTb4b4b4(
Te6edf3packetTb4b4b4: Te6edf3MeshPacketTb4b4b4,
Te6edf3pendingTb4b4b4: Te6edf3PendingResponseTb4b4b4,
Te6edf3expectedConnectionVersionTb4b4b4: Tffa657Long?Tb4b4b4,
Tb4b4b4)Tb4b4b4: Te6edf3DeferredTff7b72<Te6edf3AwaitedSendStatusTff7b72> Tb4b4b4{
Tff7b72try Tb4b4b4{
T8b949e// Publish the dispatch epoch under the response mutex so an immediate QueueStatus/Routing response cannot
T8b949e// complete this packet before transport ownership is visible.
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72val Te6edf3lifecycle Tff7b72= Te6edf3connectionStateProviderTb4b4b4.Te6edf3connectionLifecycleTb4b4b4.Te6edf3value
Tff7b72if Tb4b4b4(
Te6edf3lifecycleTb4b4b4.Te6edf3state Tff7b72!is Te6edf3ConnectionStateTb4b4b4.Te6edf3Connected Tff7b72|Tff7b72|
Tb4b4b4(Te6edf3expectedConnectionVersion Tff7b72!Tff7b72= Tff7b72null Tff7b72&Tff7b72& Te6edf3lifecycleTb4b4b4.Te6edf3version Tff7b72!Tff7b72= Te6edf3expectedConnectionVersionTb4b4b4)
Tb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3RadioNotConnectedExceptionTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tff7b72val Te6edf3departureEpoch Tff7b72= Te6edf3lifecycleTb4b4b4.Te6edf3epochsTb4b4b4.Te6edf3departures
Tff7b72val Te6edf3accepted Tff7b72= Te6edf3dispatchToRadioTb4b4b4(Te6edf3ToRadioTb4b4b4(Te6edf3packet Tff7b72= Te6edf3packetTb4b4b4)Tb4b4b4)
Te6edf3pendingTb4b4b4.Te6edf3recordDispatchTb4b4b4(Te6edf3accepted Tff7b72= Te6edf3acceptedTb4b4b4, Te6edf3departureEpoch Tff7b72= Te6edf3departureEpochTb4b4b4)
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3acceptedTb4b4b4) Te6edf3pendingTb4b4b4.Te6edf3completeAllTb4b4b4(Te6edf3AwaitedSendStatusTb4b4b4.Te6edf3SEND_FAILEDTb4b4b4)
Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3pendingTb4b4b4.Te6edf3wasDispatched Tff7b72&Tff7b72& Te6edf3pendingTb4b4b4.Te6edf3routingDeferred Tff7b72!Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{
Te6edf3scheduleRoutingResponseExpiryTb4b4b4(Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4, Te6edf3pendingTb4b4b4)
Tb4b4b4} Tff7b72else Tff7b72if Tb4b4b4(Tff7b72!Te6edf3pendingTb4b4b4.Te6edf3wasDispatchedTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffsendToRadio dropped: no active transport accepted id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff" Tb4b4b4}
Tb4b4b4}
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3exTb4b4b4: Te6edf3RadioNotConnectedExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3exTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffsendToRadio skipped: Not connected to radioTa5d6ff" Tb4b4b4}
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Te6edf3pendingTb4b4b4.Te6edf3completeAllTb4b4b4(Te6edf3AwaitedSendStatusTb4b4b4.Te6edf3TRANSPORT_STOPPEDTb4b4b4) Tb4b4b4}
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3exTb4b4b4: Te6edf3CancellationExceptionTb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3ex
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3exTb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3eTb4b4b4(Te6edf3exTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffsendToRadio error: Tffd700${Te6edf3exTb4b4b4.Te6edf3messageTffd700}Ta5d6ff" Tb4b4b4}
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Te6edf3pendingTb4b4b4.Te6edf3completeAllTb4b4b4(Te6edf3pendingTb4b4b4.Te6edf3sendFailureStatusTb4b4b4(Tb4b4b4)Tb4b4b4) Tb4b4b4}
Tb4b4b4}
Tff7b72return Te6edf3pendingTb4b4b4.Te6edf3queueDeferred
Tb4b4b4}
T8b949e/**
* Owns strict-routing expiry from service scope so caller cancellation cannot strand a packet-ID reservation. Queue
* rejection, routing completion, disconnect, and service shutdown can all finish the pending response first; the
* identity check in [removePendingIfTerminalLocked] keeps a late expiry from touching a reused ID.
*/
Tff7b72private Tff7b72fun Td2a8ffscheduleRoutingResponseExpiryTb4b4b4(Te6edf3packetIdTb4b4b4: Tffa657IntTb4b4b4, Te6edf3pendingTb4b4b4: Te6edf3PendingResponseTb4b4b4) Tb4b4b4{
Te6edf3scopeTb4b4b4.Te6edf3handledLaunch Tb4b4b4{
Te6edf3delayTb4b4b4(Te6edf3ROUTING_RESPONSE_TIMEOUTTb4b4b4)
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3pendingTb4b4b4.Te6edf3completeRoutingTb4b4b4(Te6edf3AwaitedSendStatusTb4b4b4.Te6edf3TIMED_OUTTb4b4b4)
Te6edf3removePendingIfTerminalLockedTb4b4b4(Te6edf3packetIdTb4b4b4, Te6edf3pendingTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffremovePendingIfTerminalTb4b4b4(Te6edf3packetIdTb4b4b4: Tffa657IntTb4b4b4, Te6edf3pendingTb4b4b4: Te6edf3PendingResponseTb4b4b4) Tff7b72=
Te6edf3withContextTb4b4b4(Te6edf3NonCancellableTb4b4b4) Tb4b4b4{ Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Te6edf3removePendingIfTerminalLockedTb4b4b4(Te6edf3packetIdTb4b4b4, Te6edf3pendingTb4b4b4) Tb4b4b4} Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffremovePendingIfTerminalLockedTb4b4b4(Te6edf3packetIdTb4b4b4: Tffa657IntTb4b4b4, Te6edf3pendingTb4b4b4: Te6edf3PendingResponseTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3pendingTb4b4b4.Te6edf3isTerminal Tff7b72&Tff7b72& Te6edf3queueResponseTff7b72[Te6edf3packetIdTff7b72] Tff7b72=Tff7b72=Tff7b72= Te6edf3pendingTb4b4b4) Te6edf3queueResponseTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3packetIdTb4b4b4)
Tb4b4b4}
Tff7b72private Tff7b72fun Te6edf3BooleanTb4b4b4.Td2a8fftoAwaitedSendStatusTb4b4b4(Tb4b4b4)Tb4b4b4: Te6edf3AwaitedSendStatus Tff7b72=
Tff7b72if Tb4b4b4(Tff7b72thisTb4b4b4) Te6edf3AwaitedSendStatusTb4b4b4.Te6edf3ACCEPTED Tff7b72else Te6edf3AwaitedSendStatusTb4b4b4.Te6edf3RADIO_REJECTED
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffcompletePendingResponsesTb4b4b4(Te6edf3statusTb4b4b4: Te6edf3AwaitedSendStatusTb4b4b4)Tb4b4b4: Te6edf3ListTff7b72<Te6edf3PacketStatusTargetTff7b72> Tff7b72=
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72val Te6edf3completedPackets Tff7b72=
Te6edf3queueResponseTb4b4b4.Te6edf3mapNotNull Tb4b4b4{ Tb4b4b4(Te6edf3_Tb4b4b4, Te6edf3pendingTb4b4b4) Tff7b72-Tff7b72>
Te6edf3pendingTb4b4b4.Te6edf3completeAllTb4b4b4(Te6edf3statusTb4b4b4)
Te6edf3PacketStatusTargetTb4b4b4(Te6edf3pendingTb4b4b4.Te6edf3packetTb4b4b4, Te6edf3pendingTb4b4b4.Te6edf3persistenceTb4b4b4)Tb4b4b4.Te6edf3takeIf Tb4b4b4{
Te6edf3pendingTb4b4b4.Te6edf3terminalStatus Tff7b72!Tff7b72= Te6edf3AwaitedSendStatusTb4b4b4.Te6edf3ACCEPTED
Tb4b4b4}
Tb4b4b4}
Te6edf3queueResponseTb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)
Te6edf3completedPackets
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffinsertMeshLogTb4b4b4(Te6edf3packetToSaveTb4b4b4: Te6edf3MeshLogTb4b4b4) Tb4b4b4{
Te6edf3scopeTb4b4b4.Te6edf3handledLaunch Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{
Ta5d6ff"Ta5d6ffinsert: Tffd700${Te6edf3packetToSaveTb4b4b4.Te6edf3message_typeTffd700}Ta5d6ff = Ta5d6ff" Tff7b72+
Ta5d6ff"Tffd700${Te6edf3packetToSaveTb4b4b4.Te6edf3raw_messageTb4b4b4.Te6edf3toOneLineStringTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff from=Tffd700${Te6edf3packetToSaveTb4b4b4.Te6edf3fromNumTffd700}Ta5d6ff"
Tb4b4b4}
Te6edf3meshLogRepositoryTb4b4b4.Te6edf3valueTb4b4b4.Te6edf3insertTb4b4b4(Te6edf3packetToSaveTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Served by rngit 1.5.0 - Generated in 0.42s